Repository navigation
fix(database): restack Mongo live updates onto current Event Relays - #1606
Merged
kkopanidis merged 6 commits intoSep 13, 2026
Merged
Conversation
…ger) (#1604) * fix(router): harden event relays for HA and auth lifetime EventBus uses a single Redis message dispatcher with subscriberId maps, SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine middleware once, skips empty-room namespace broadcasts, demotes hot-path logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager gains coalesced periodic reconcile, compiled templates, inbound caps, TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and metrics/docs aligned with the interventions plan. * fix(router): address #1604 re-review P1 and P2 follow-ups Recovery re-auth keeps per-user subscription state across disconnect; emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process signal handlers and adds subscribeAck; manager subscribes only after Redis ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local sockets without emit acks; cache is bounded; preview caps sample JSON; handshake matcher rejects ticket/sid/non-namespace paths; socket middleware rebind avoids duplicate registration on first sockets enable. * fix(hermes): restore production tsc for engine middleware chain Type the Socket.IO engine middleware runner as Express NextFunction so recursive callbacks match registerGlobalMiddleware. Use Logger.info for socket trace lines (IConduitLogger has no debug). * fix(router): import Express Response for socket global middleware _rebindSocketGlobalMiddlewares passes handlers typed with fetch Response because Response was not imported from express, failing tsc against Hermes registerSocketGlobalMiddleware. * fix(router): quit EventBus on module shutdown signals Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(), registered on SIGTERM/SIGINT with process exit so Redis teardown is not left to SDK signal handlers. Tighten registerGlobalMiddleware typing and trim handshake helper comment (deslop). * fix(router): pass-3 recovery, scoped emit, and subscription hygiene Recovery re-auth keeps membership on Authorization UNAVAILABLE and only leaves on deny; re-check subscriptions for restored er: rooms via room map. Prune user subscriptions on disconnect when no other socket holds them. Scope relay emits to receivers in the target room; use Engine.IO writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs. * fix(router): pass-4 recovery map, backpressure, emit tests Store relay subs on socket.data; resolve recovered rooms from context and TTL room map without last-writer userId. Prune room map when unused on disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router tests for room-scoped emit and queue metric. * fix(router): clear eventRelaySubs from socket.data on recovery deny * fix(router): clear socket.data subs on recovery fail-closed leave
Split validateEventRelayInput into per-field parsers to reduce complexity. Silence unused-parameter lint in EventBus test FakeRedis stub.
Replay the Mongo delta onto feat/router-event-relays (9e6356e). Keep #1604 Hermes (engine.use, er: recovery, writeBuffer) and graft only the port-keyed adapter plus register-before-initSockets. Do not take old #1601 Socket.ts wholesale. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com>
…is keys Lock renew failure bumps a generation so a draining watch cannot emit. Do not open a watch unless the leader lock extends after acquire. Arm TTL on realtime socket/doc keys when a recoverable disconnect never recovers; persist on restore. CMS-read deny fails closed at emit without extra ReBAC round-trips and keeps membership; UNAVAILABLE keeps membership. Cover dropDatabase reopen. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com>
… watches Two overlapping ensureLeader calls both bumped the lock generation while the first openStream was still reading the resume token, so the watch never attached. Hold an acquiring flag so only one opener runs.
kkopanidis
marked this pull request as ready for review
September 13, 2026 18:54
2 of 3 tasks
kkopanidis
added a commit
that referenced
this pull request
Oct 6, 2026
* feat(database): stream live document updates over admin sockets MongoDB change streams require a replica set, so local Compose now initializes a single-node rs0. Schemas can opt in via realtime.enabled. * fix(docker): initiate local Mongo replica set before Conduit starts Mongo 4.4 keyfiles reject hyphens, and rs.status() returns ok:0 instead of throwing, so Compose never elected a primary for replicaSet URIs. * fix: isolate socket streams and persist realtime schema options Stop duplicate admin/router deliveries and apply handshake auth on connect. Retry change-stream leadership after lock races. * refactor(database): drop unused realtime aliases and tautological tests * fix(database): Mongo live-update interventions (watch, token, tickets, recovery) (#1605) * fix(router): Event Relays interventions (EventBus, Hermes, relay manager) (#1604) * fix(router): harden event relays for HA and auth lifetime EventBus uses a single Redis message dispatcher with subscriberId maps, SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine middleware once, skips empty-room namespace broadcasts, demotes hot-path logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager gains coalesced periodic reconcile, compiled templates, inbound caps, TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and metrics/docs aligned with the interventions plan. * fix(router): address #1604 re-review P1 and P2 follow-ups Recovery re-auth keeps per-user subscription state across disconnect; emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process signal handlers and adds subscribeAck; manager subscribes only after Redis ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local sockets without emit acks; cache is bounded; preview caps sample JSON; handshake matcher rejects ticket/sid/non-namespace paths; socket middleware rebind avoids duplicate registration on first sockets enable. * fix(hermes): restore production tsc for engine middleware chain Type the Socket.IO engine middleware runner as Express NextFunction so recursive callbacks match registerGlobalMiddleware. Use Logger.info for socket trace lines (IConduitLogger has no debug). * fix(router): import Express Response for socket global middleware _rebindSocketGlobalMiddlewares passes handlers typed with fetch Response because Response was not imported from express, failing tsc against Hermes registerSocketGlobalMiddleware. * fix(router): quit EventBus on module shutdown signals Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(), registered on SIGTERM/SIGINT with process exit so Redis teardown is not left to SDK signal handlers. Tighten registerGlobalMiddleware typing and trim handshake helper comment (deslop). * fix(router): pass-3 recovery, scoped emit, and subscription hygiene Recovery re-auth keeps membership on Authorization UNAVAILABLE and only leaves on deny; re-check subscriptions for restored er: rooms via room map. Prune user subscriptions on disconnect when no other socket holds them. Scope relay emits to receivers in the target room; use Engine.IO writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs. * fix(router): pass-4 recovery map, backpressure, emit tests Store relay subs on socket.data; resolve recovered rooms from context and TTL room map without last-writer userId. Prune room map when unused on disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router tests for room-scoped emit and queue metric. * fix(router): clear eventRelaySubs from socket.data on recovery deny * fix(router): clear socket.data subs on recovery fail-closed leave * fix(router): address CodeFactor findings on event relay validation Split validateEventRelayInput into per-field parsers to reduce complexity. Silence unused-parameter lint in EventBus test FakeRedis stub. * fix(database): Mongo live-update interventions (watch, token, tickets, recovery) Rebase #1601 onto current Event Relays Hermes and close the stack-review P1s: watch $match/$project, serialized handleChange, persist resume token after successful emit, keep token on CursorKilled 237, hermes isSocketHandshake, restore /database/ Redis on recover, and shutdown the watch plus leader lock. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> * style: prettier on Mongo intervention tests * fix(database): stop watch on emit failure; keep ReBAC membership Failed socketPush/bus no longer advances the resume token. ReBAC unavailable no longer removeUser. Auth middleware 401s a realtime ticket on POST /realtime/ticket. Drop/rename/invalidate reopen the watch. --------- * fix(database): restack Mongo live updates onto current Event Relays (#1606) * fix(router): Event Relays interventions (EventBus, Hermes, relay manager) (#1604) * fix(router): harden event relays for HA and auth lifetime EventBus uses a single Redis message dispatcher with subscriberId maps, SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine middleware once, skips empty-room namespace broadcasts, demotes hot-path logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager gains coalesced periodic reconcile, compiled templates, inbound caps, TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and metrics/docs aligned with the interventions plan. * fix(router): address #1604 re-review P1 and P2 follow-ups Recovery re-auth keeps per-user subscription state across disconnect; emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process signal handlers and adds subscribeAck; manager subscribes only after Redis ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local sockets without emit acks; cache is bounded; preview caps sample JSON; handshake matcher rejects ticket/sid/non-namespace paths; socket middleware rebind avoids duplicate registration on first sockets enable. * fix(hermes): restore production tsc for engine middleware chain Type the Socket.IO engine middleware runner as Express NextFunction so recursive callbacks match registerGlobalMiddleware. Use Logger.info for socket trace lines (IConduitLogger has no debug). * fix(router): import Express Response for socket global middleware _rebindSocketGlobalMiddlewares passes handlers typed with fetch Response because Response was not imported from express, failing tsc against Hermes registerSocketGlobalMiddleware. * fix(router): quit EventBus on module shutdown signals Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(), registered on SIGTERM/SIGINT with process exit so Redis teardown is not left to SDK signal handlers. Tighten registerGlobalMiddleware typing and trim handshake helper comment (deslop). * fix(router): pass-3 recovery, scoped emit, and subscription hygiene Recovery re-auth keeps membership on Authorization UNAVAILABLE and only leaves on deny; re-check subscriptions for restored er: rooms via room map. Prune user subscriptions on disconnect when no other socket holds them. Scope relay emits to receivers in the target room; use Engine.IO writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs. * fix(router): pass-4 recovery map, backpressure, emit tests Store relay subs on socket.data; resolve recovered rooms from context and TTL room map without last-writer userId. Prune room map when unused on disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router tests for room-scoped emit and queue metric. * fix(router): clear eventRelaySubs from socket.data on recovery deny * fix(router): clear socket.data subs on recovery fail-closed leave * fix(router): address CodeFactor findings on event relay validation Split validateEventRelayInput into per-field parsers to reduce complexity. Silence unused-parameter lint in EventBus test FakeRedis stub. * feat(database): restack Mongo live updates onto current Event Relays Replay the Mongo delta onto feat/router-event-relays (9e6356e). Keep #1604 Hermes (engine.use, er: recovery, writeBuffer) and graft only the port-keyed adapter plus register-before-initSockets. Do not take old #1601 Socket.ts wholesale. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> * fix(database): fence change-stream leadership and expire recovery Redis keys Lock renew failure bumps a generation so a draining watch cannot emit. Do not open a watch unless the leader lock extends after acquire. Arm TTL on realtime socket/doc keys when a recoverable disconnect never recovers; persist on restore. CMS-read deny fails closed at emit without extra ReBAC round-trips and keeps membership; UNAVAILABLE keeps membership. Cover dropDatabase reopen. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> * fix(database): serialize leader acquire so concurrent reconcile still watches Two overlapping ensureLeader calls both bumped the lock generation while the first openStream was still reading the resume token, so the watch never attached. Hold an acquiring flag so only one opener runs. --------- * fix(database): watch Mongo live updates from now, drop resume tokens Live updates notify only: change streams start at the end of the oplog. Do not persist Redis resume tokens or restore on CursorKilled. Leader restart or cursor drop is a gap; clients refetch. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> * fix(database): drop leftover resumeAfter watch options Watch is pipeline-only (db.watch(pipeline)). Strip unused change-stream _id/clusterTime fields and test resume-token fixtures. Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> * fix(database): fence Mongo change-stream leader and reset on new watch (#1625) Stop a frozen leader from emitting after takeover by INCR-ing a Redis fence on lock acquire and re-reading it before every database:change publish and socket push. Poll lock contention about every second, keep exponential backoff for stream errors, wake standbys on clean lock release, and push one schema-room reset when a watch opens. Co-authored-by: Cursor Agent <cursoragent@cursor.com> --------- Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com> Co-authored-by: Cursor Agent <cursoragent@cursor.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Restacks #1601 onto current #1600 (
9e6356ef) and lands leftover Mongo interventions. Head is strictly ahead of #1600. Do not merge #1600 or #1601. Do not retarget tomain.Restack
feat/router-event-relays(9e6356ef). Did not start from squashbc323840.engine.use,er:recovery,writeBuffer). Graft only port-keyed adaptersocket.io:${port}and register-before-initSockets.engine.useandapplySocketGlobalMiddlewares.Leftovers
realtime:socket:*/ exclusiverealtime:doc:*keys expire after 120s if recovery never happens. Restore persists keys.dropDatabasereopen test alongside drop/rename/invalidate.$match/$project; skip path is token-only; serializedhandleChange; batched ReBAC+cache. Token persists after successful emit. NofullDocumentin the payload.Locked
stopStream+scheduleRetry.router.eventRelays.enabled, no tenants.Out of scope
getLocalRoomUserIds/getLocalRoomsWithPrefix